welo git mirrors
src/mux.rs 5fdea09c4ee193f99cd5512f1a8d81c363b4012e (5fdea09c) Text, 7.49 KB
T8b949e//! Multiplexer -- manages sessions over a single RNS link.
T8b949e//!
T8b949e//! All TCP sessions are multiplexed over one RNS link per peer. Each session
T8b949e//! gets a unique `session_id`. Data is exchanged as frames:
T8b949e//!
T8b949e//! ```text
T8b949e//! [1 byte type][4 bytes session_id][2 bytes payload length][payload]
T8b949e//! ```
T8b949e//!
T8b949e//! Frames are sent as raw encrypted link data via `send_on_link` (CONTEXT_NONE).
T8b949e//! Frames larger than LINK_MDU are split into LINK_MDU-sized chunks; the
T8b949e//! receiver reassembles them using the frame's embedded length header.
T8b949e//!
T8b949e//! We do NOT use the RNS Channel API because its ACK path is not wired up
T8b949e//! in rns-rs, causing `NotReady` errors after the first 2 messages.
Tff7b72use Te6edf3std::Te6edf3collections::Te6edf3HashMapTb4b4b4;
Tff7b72use Te6edf3std::Te6edf3errorTb4b4b4;
Tff7b72use Te6edf3std::Te6edf3sync::Tb4b4b4{Te6edf3ArcTb4b4b4, Te6edf3MutexTb4b4b4}Tb4b4b4;
Tff7b72use Te6edf3log::Tb4b4b4{Te6edf3debugTb4b4b4, Te6edf3errorTb4b4b4, Te6edf3infoTb4b4b4, Te6edf3warnTb4b4b4}Tb4b4b4;
Tff7b72use Te6edf3rns_core::Te6edf3constants::Te6edf3LINK_MDUTb4b4b4;
Tff7b72use Te6edf3rns_core::Te6edf3msgpack::Te6edf3ErrorTb4b4b4;
Tff7b72use Te6edf3rns_net::Tb4b4b4{Te6edf3LinkIdTb4b4b4, Te6edf3LocalServerFactoryTb4b4b4, Te6edf3RnsNodeTb4b4b4}Tb4b4b4;
Tff7b72use Te6edf3tokio::Te6edf3sync::Te6edf3mpscTb4b4b4;
Tff7b72use Tff7b72crate::Te6edf3frame::Te6edf3FrameDecodeState::Tb4b4b4{Te6edf3DecodingFailedTb4b4b4, Te6edf3MoreDataRequiredTb4b4b4}Tb4b4b4;
Tff7b72use Tff7b72crate::Tb4b4b4{Te6edf3FrameTb4b4b4, Te6edf3FrameTypeTb4b4b4}Tb4b4b4;
T8b949e/// Context byte for our link data. We use CONTEXT_NONE (0x00) which routes
T8b949e/// to the `on_link_data` callback on the receive side.
Tff7b72const Te6edf3DATA_CONTEXT: Tffa657u8 Tff7b72= T79c0ff0x00Tb4b4b4;
T8b949e/// A handle that session tasks use to send frames back through the RNS link.
Tff7b72#[Tff7b72derive(Clone)Tff7b72]
Tff7b72pub Tff7b72struct T56d364MuxHandle Tb4b4b4{
Te6edf3inner: T56d364ArcTff7b72<Te6edf3MuxInnerTff7b72>Tb4b4b4,
Tb4b4b4}
Tff7b72struct T56d364MuxInner Tb4b4b4{
Te6edf3node: T56d364ArcTff7b72<Te6edf3RnsNodeTff7b72>Tb4b4b4,
Te6edf3link_id: T56d364MutexTff7b72<Tffa657OptionTff7b72<Te6edf3LinkIdTff7b72>Tff7b72>Tb4b4b4,
Te6edf3sessions: T56d364MutexTff7b72<Te6edf3HashMapTff7b72<Tffa657u32Tb4b4b4, Te6edf3mpsc::Te6edf3UnboundedSenderTff7b72<Te6edf3FrameTff7b72>Tff7b72>Tff7b72>Tb4b4b4,
Te6edf3next_sid: T56d364MutexTff7b72<Tffa657u32Tff7b72>Tb4b4b4,
T8b949e/// Reassembly buffer for incoming raw link data chunks.
Te6edf3recv_buf: T56d364MutexTff7b72<Tffa657VecTff7b72<Tffa657u8Tff7b72>Tff7b72>Tb4b4b4,
Tb4b4b4}
Tff7b72impl Te6edf3MuxHandle Tb4b4b4{
T8b949e/// Create a new multiplexer handle.
Tff7b72pub Tff7b72fn Td2a8ffnewTb4b4b4(Te6edf3node: T56d364ArcTff7b72<Te6edf3RnsNodeTff7b72>Tb4b4b4) -> T56d364Self Tb4b4b4{
Tff7b72Self Tb4b4b4{
Te6edf3inner: T56d364Arc::Te6edf3newTb4b4b4(Te6edf3MuxInner Tb4b4b4{
Te6edf3nodeTb4b4b4,
Te6edf3link_id: T56d364Mutex::Te6edf3newTb4b4b4(Tffa657NoneTb4b4b4)Tb4b4b4,
Te6edf3sessions: T56d364Mutex::Te6edf3newTb4b4b4(Te6edf3HashMap::Te6edf3newTb4b4b4(Tb4b4b4)Tb4b4b4)Tb4b4b4,
Te6edf3next_sid: T56d364Mutex::Te6edf3newTb4b4b4(T79c0ff0Tb4b4b4)Tb4b4b4,
Te6edf3recv_buf: T56d364Mutex::Te6edf3newTb4b4b4(Tffa657Vec::Te6edf3newTb4b4b4(Tb4b4b4)Tb4b4b4)Tb4b4b4,
Tb4b4b4}Tb4b4b4)Tb4b4b4,
Tb4b4b4}
Tb4b4b4}
T8b949e/// Set the active link id (called when the link is established).
Tff7b72pub Tff7b72fn Td2a8ffset_link_idTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4, Te6edf3link_id: T56d364LinkIdTb4b4b4) Tb4b4b4{
Tff7b72*Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3link_idTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4) Tff7b72= Tffa657SomeTb4b4b4(Te6edf3link_idTb4b4b4)Tb4b4b4;
Tb4b4b4}
T8b949e/// Clear the link id (called when the link is closed).
Tff7b72pub Tff7b72fn Td2a8ffclear_link_idTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4) Tb4b4b4{
Tff7b72*Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3link_idTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4) Tff7b72= Tffa657NoneTb4b4b4;
Tb4b4b4}
T8b949e/// Reset the mux for a new link (clears sessions and reassembly buffer).
T8b949e///
T8b949e/// Called on reconnection.
Tff7b72pub Tff7b72fn Td2a8ffresetTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4) Tb4b4b4{
Tff7b72*Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3link_idTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4) Tff7b72= Tffa657NoneTb4b4b4;
Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3sessionsTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)Tb4b4b4;
Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3recv_bufTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)Tb4b4b4;
Tb4b4b4}
T8b949e/// Check if link is active.
Tff7b72pub Tff7b72fn Td2a8ffis_connectedTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4) -> Tffa657bool Tb4b4b4{
Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3link_idTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3is_someTb4b4b4(Tb4b4b4)
Tb4b4b4}
T8b949e/// Get the next session id.
Tff7b72pub Tff7b72fn Td2a8ffnext_session_idTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4) -> Tffa657u32 Tb4b4b4{
Tff7b72let Tff7b72mut Te6edf3sid Tff7b72= Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3next_sidTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4;
Tff7b72*Te6edf3sid Tff7b72= Te6edf3sidTb4b4b4.Te6edf3wrapping_addTb4b4b4(T79c0ff1Tb4b4b4)Tb4b4b4;
Tff7b72*Te6edf3sid
Tb4b4b4}
T8b949e/// Register a session. Returns a receiver for frames addressed to this session.
Tff7b72pub Tff7b72fn Td2a8ffregister_sessionTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4, Te6edf3sid: Tffa657u32Tb4b4b4) -> T56d364mpsc::Te6edf3UnboundedReceiverTff7b72<Te6edf3FrameTff7b72> Tb4b4b4{
Tff7b72let Tb4b4b4(Te6edf3txTb4b4b4, Te6edf3rxTb4b4b4) Tff7b72= Te6edf3mpsc::Te6edf3unbounded_channelTb4b4b4(Tb4b4b4)Tb4b4b4;
Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3sessionsTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3insertTb4b4b4(Te6edf3sidTb4b4b4, Te6edf3txTb4b4b4)Tb4b4b4;
Te6edf3rx
Tb4b4b4}
T8b949e/// Remove a session.
Tff7b72pub Tff7b72fn Td2a8ffdrop_sessionTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4, Te6edf3sid: Tffa657u32Tb4b4b4) Tb4b4b4{
Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3sessionsTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3removeTb4b4b4(Tff7b72&Te6edf3sidTb4b4b4)Tb4b4b4;
Tb4b4b4}
T8b949e/// Send a frame over the RNS link.
T8b949e///
T8b949e/// The encoded frame is split into chunks of at most `LINK_MDU` bytes and
T8b949e/// each chunk is sent as a separate `send_on_link` call with `CONTEXT_NONE`.
T8b949e/// The receiver reassembles using the frame length header.
T8b949e///
T8b949e/// Note that we are using .inner.link_id.lock as the guard to prevent
T8b949e/// multiple different sids from sending at the same time and scrambling packets
Tff7b72pub Tff7b72fn Td2a8ffsend_frameTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4, Te6edf3frame: Tff7b72&T56d364FrameTb4b4b4) Tb4b4b4{
Tff7b72let Te6edf3lock Tff7b72= Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3link_idTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4;
Tff7b72let Te6edf3value Tff7b72= Tff7b72&Tff7b72*Tb4b4b4(Te6edf3lockTb4b4b4)Tb4b4b4;
Tff7b72let Te6edf3link_id Tff7b72= Tff7b72match Te6edf3value Tb4b4b4{
Tffa657SomeTb4b4b4(Te6edf3idTb4b4b4) Tff7b72=Tff7b72> Te6edf3idTb4b4b4,
Tffa657None Tff7b72=Tff7b72> Tb4b4b4{
Te6edf3debugTff7b72!Tb4b4b4(Ta5d6ff"Ta5d6ffsend_frame: no active link, dropping frameTa5d6ff"Tb4b4b4)Tb4b4b4;
Tff7b72returnTb4b4b4;
Tb4b4b4}
Tb4b4b4}Tb4b4b4;
Tff7b72let Te6edf3encoded Tff7b72= Te6edf3frameTb4b4b4.Te6edf3encodeTb4b4b4(Tb4b4b4)Tb4b4b4;
Tff7b72let Te6edf3node Tff7b72= Tff7b72&Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3nodeTb4b4b4;
T8b949e// Send in LINK_MDU-sized chunks
Tff7b72for Te6edf3chunk Tff7b72in Te6edf3encodedTb4b4b4.Te6edf3chunksTb4b4b4(Te6edf3LINK_MDUTb4b4b4) Tb4b4b4{
T8b949e// println!("chunk: {:?} sid: {}",chunk, frame.session_id);
Tff7b72if Tff7b72let Tffa657ErrTb4b4b4(Te6edf3eTb4b4b4) Tff7b72= Te6edf3nodeTb4b4b4.Te6edf3send_on_linkTb4b4b4(Te6edf3link_idTb4b4b4.T79c0ff0Tb4b4b4, Te6edf3chunkTb4b4b4.Te6edf3to_vecTb4b4b4(Tb4b4b4)Tb4b4b4, Te6edf3DATA_CONTEXTTb4b4b4) Tb4b4b4{
Te6edf3warnTff7b72!Tb4b4b4(Ta5d6ff"Ta5d6ffFailed to send link data: {:?}Ta5d6ff"Tb4b4b4, Te6edf3eTb4b4b4)Tb4b4b4;
Tff7b72returnTb4b4b4;
Tb4b4b4}
Tb4b4b4}
Tffa657dropTb4b4b4(Te6edf3lockTb4b4b4)Tb4b4b4;
Tb4b4b4}
T8b949e/// Convenience: send a typed frame.
Tff7b72pub Tff7b72fn Td2a8ffsendTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4, Te6edf3frame_type: T56d364FrameTypeTb4b4b4, Te6edf3session_id: Tffa657u32Tb4b4b4, Te6edf3payload: Tffa657VecTff7b72<Tffa657u8Tff7b72>Tb4b4b4) Tb4b4b4{
T8b949e// info!("send frame");
Tff7b72selfTb4b4b4.Te6edf3send_frameTb4b4b4(Tff7b72&Te6edf3Frame::Te6edf3newTb4b4b4(Te6edf3frame_typeTb4b4b4, Te6edf3session_idTb4b4b4, Te6edf3payloadTb4b4b4)Tb4b4b4)Tb4b4b4;
Tb4b4b4}
T8b949e/// Dispatch an incoming frame to the appropriate session.
Tff7b72pub Tff7b72fn Td2a8ffdispatchTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4, Te6edf3frame: T56d364FrameTb4b4b4) Tb4b4b4{
Tff7b72let Te6edf3sid Tff7b72= Te6edf3frameTb4b4b4.Te6edf3session_idTb4b4b4;
Tff7b72let Te6edf3ft Tff7b72= Te6edf3frameTb4b4b4.Te6edf3frame_typeTb4b4b4;
Tff7b72let Te6edf3sessions Tff7b72= Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3sessionsTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4;
Tff7b72if Tff7b72let Tffa657SomeTb4b4b4(Te6edf3txTb4b4b4) Tff7b72= Te6edf3sessionsTb4b4b4.Te6edf3getTb4b4b4(Tff7b72&Te6edf3sidTb4b4b4) Tb4b4b4{
Tff7b72if Te6edf3txTb4b4b4.Te6edf3sendTb4b4b4(Te6edf3frameTb4b4b4)Tb4b4b4.Te6edf3is_errTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3debugTff7b72!Tb4b4b4(Ta5d6ff"Ta5d6ffSession {} channel closedTa5d6ff"Tb4b4b4, Te6edf3sidTb4b4b4)Tb4b4b4;
Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3warnTff7b72!Tb4b4b4(Ta5d6ff"Ta5d6ffNo session {} for frame type {}Ta5d6ff"Tb4b4b4, Te6edf3sidTb4b4b4, Te6edf3ftTb4b4b4)Tb4b4b4;
Tb4b4b4}
Tb4b4b4}
T8b949e/// Feed raw link data and extract any complete frames.
T8b949e///
T8b949e/// Called when `on_link_data` fires. The data may be a partial chunk of a
T8b949e/// larger frame, so we buffer and try to decode complete frames.
Tff7b72pub Tff7b72fn Td2a8ffreceive_dataTb4b4b4(Tff7b72&Tff7b72selfTb4b4b4, Te6edf3data: Tff7b72&Tb4b4b4[Tffa657u8Tb4b4b4]Tb4b4b4) -> Tffa657VecTff7b72<Te6edf3FrameTff7b72> Tb4b4b4{
Tff7b72let Tff7b72mut Te6edf3buf Tff7b72= Tff7b72selfTb4b4b4.Te6edf3innerTb4b4b4.Te6edf3recv_bufTb4b4b4.Te6edf3lockTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3unwrapTb4b4b4(Tb4b4b4)Tb4b4b4;
Te6edf3bufTb4b4b4.Te6edf3extend_from_sliceTb4b4b4(Te6edf3dataTb4b4b4)Tb4b4b4;
T8b949e// info!("buf: {:?}", buf);
Tff7b72let Tff7b72mut Te6edf3frames Tff7b72= Tffa657Vec::Te6edf3newTb4b4b4(Tb4b4b4)Tb4b4b4;
Tff7b72loop Tb4b4b4{
Tff7b72match Te6edf3Frame::Te6edf3decodeTb4b4b4(Tff7b72&Te6edf3bufTb4b4b4) Tb4b4b4{
Tffa657OkTb4b4b4(Tb4b4b4(Te6edf3frameTb4b4b4, Te6edf3consumedTb4b4b4)Tb4b4b4) Tff7b72=Tff7b72> Tb4b4b4{
T8b949e// info!("consume {}",consumed);
Te6edf3bufTb4b4b4.Te6edf3drainTb4b4b4(Tff7b72..Te6edf3consumedTb4b4b4)Tb4b4b4;
Te6edf3framesTb4b4b4.Te6edf3pushTb4b4b4(Te6edf3frameTb4b4b4)Tb4b4b4;
Tb4b4b4}
Tffa657ErrTb4b4b4(Te6edf3errTb4b4b4) Tff7b72=Tff7b72> Tb4b4b4{
Tff7b72match Te6edf3err Tb4b4b4{
T8b949e// Finished => {
T8b949e// // finished decoding all packets basically
T8b949e// // though there might still be like 5 bytes left in the buffer
T8b949e// info!("finished");
T8b949e// break;
T8b949e// }
Te6edf3MoreDataRequired Tff7b72=Tff7b72> Tb4b4b4{
Te6edf3infoTff7b72!Tb4b4b4(Ta5d6ff"Ta5d6ffmore data requiredTa5d6ff"Tb4b4b4)Tb4b4b4;
Tff7b72breakTb4b4b4;
T8b949e// just wait for next packet
Tb4b4b4}
Te6edf3DecodingFailed Tff7b72=Tff7b72> Tb4b4b4{
T8b949e// something has gone really wrong, just clear the buffer and hope things
T8b949e// work out.
Te6edf3errorTff7b72!Tb4b4b4(Ta5d6ff"Ta5d6ffdecoding failed for a packet, something really bad is happening {:?}Ta5d6ff"Tb4b4b4,Te6edf3bufTb4b4b4)Tb4b4b4;
Te6edf3errorTff7b72!Tb4b4b4(Ta5d6ff"Ta5d6ffignoring packet buffer and hoping that will fix itTa5d6ff"Tb4b4b4)Tb4b4b4;
Tff7b72assert!Tb4b4b4(Tff7b72falseTb4b4b4)Tb4b4b4;
Te6edf3bufTb4b4b4.Te6edf3drainTb4b4b4(Tff7b72..Tb4b4b4)Tb4b4b4;
Tff7b72breakTb4b4b4;
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}Tb4b4b4,
Tb4b4b4}
Tb4b4b4}
T8b949e// // info!("frames {:?}", &frames);
T8b949e// if buf.len() > 12000 {
T8b949e// warn!("buf thing: {:?}", &buf);
T8b949e// // assert!(false);
T8b949e// }
Te6edf3frames
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.1 - Generated in 0.04s